跳到主要内容

指定时间偏移量

指定时间消费

public class KafkaOffsetSearch {
public static void main(String[] args) {
DateFormat df = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
Date now = new Date();
long nowTime = now.getTime();
System.out.println("当前时间: " + df.format(now));
// long fetchDataTime = nowTime - 1000 * 60 * 1820; // 计算一天前的时间戳
long fetchDataTime = nowTime - 1000 * 60 * 60 * 24 * 2; // 计算一天前的时间戳
String serverIp = "ip:9092";
String groupId = "test";
String topics="topic1,topic2";
String[] singleTopic = topics.split(",");
for (String topicItem : singleTopic) {
restartOffsetFromTime(serverIp, topicItem, groupId, fetchDataTime);
}
}

public static void restartOffsetFromTime(String serverIp, String topic, String groupId, long startTime) {

Properties props = new Properties();
props.put("bootstrap.servers", serverIp);
props.put("group.id", groupId);
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
try {
// 获取topic的partition信息
List<PartitionInfo> partitionInfos = consumer.partitionsFor(topic);
List<TopicPartition> topicPartitions = new ArrayList<>();

Map<TopicPartition, Long> timestampsToSearch = new HashMap<>();


for (PartitionInfo partitionInfo : partitionInfos) {
topicPartitions.add(new TopicPartition(partitionInfo.topic(), partitionInfo.partition()));
timestampsToSearch.put(new TopicPartition(partitionInfo.topic(), partitionInfo.partition()), startTime);
}

consumer.assign(topicPartitions);

// 获取每个partition一个小时之前的偏移量,这里得到的是写入kafka消息的数据
Map<TopicPartition, OffsetAndTimestamp> map = consumer.offsetsForTimes(timestampsToSearch);
DateFormat df = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
OffsetAndTimestamp offsetTimestamp = null;
System.out.println("开始设置各分区初始偏移量...");
for (Map.Entry<TopicPartition, OffsetAndTimestamp> entry : map.entrySet()) {
// 如果设置的查询偏移量的时间点大于最大的索引记录时间,那么value就为空
offsetTimestamp = entry.getValue();
if (offsetTimestamp != null) {
int partition = entry.getKey().partition();
long timestamp = offsetTimestamp.timestamp();
long offset = offsetTimestamp.offset();
System.out.println("partition = " + partition +
", time = " + df.format(new Date(timestamp)) +
", offset = " + offset);
// 设置读取消息的偏移量
consumer.seek(entry.getKey(), offset);
}
}
System.out.println("设置各分区初始偏移量结束...");

// while(true) {
// ConsumerRecords<String, String> records = consumer.poll(1000);
// for (ConsumerRecord<String, String> record : records) {
// System.out.println("partition = " + record.partition() + ", offset = " + record.offset());
// }
// }
} catch (Exception e) {
e.printStackTrace();
} finally {
consumer.close();
}
}

public static Set<String> getAllTopicFromGroupId(String serverIp, String groupId) {
Properties props = new Properties();
props.put("bootstrap.servers", serverIp);
props.put("group.id", "test");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
return consumer.listTopics().keySet();
}
}